package SinkTest

import org.apache.flink.api.common.serialization.SimpleStringSchema
import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer011

object KafkaSink {
  def main(args: Array[String]): Unit = {

    val env = StreamExecutionEnvironment.getExecutionEnvironment
    val stream1 = env.readTextFile("path")

    stream1.addSink(
      new FlinkKafkaProducer011[String](
        "localhost:9092", "test", new SimpleStringSchema()
      )
    )

    env.execute()
  }

}
